Repository navigation
[fix][client] Return flow-control permit when client seek by messageID - #26726
programmerahul wants to merge 6 commits into
Conversation
…geId boundary message When a consumer seeks to (or is created with) a specific startMessageId, the broker re-dispatches the boundary message, which the client filters out in messageReceived() via the `isSameEntry(msgId) && isPriorEntryIndex(...)` block. That block released the payload and returned without calling increaseAvailablePermits(), so the dropped message's outstanding flow-control permit was never repaid (it is not delivered, so messageProcessed() never runs for it). Each seek-to-messageId therefore leaks one permit. It is masked at normal receiver queue sizes (seek re-establishes flow control), but with receiverQueueSize=1 the single leaked permit exhausts the whole budget and the consumer stalls immediately after the seek -- the message after the seek target is never delivered. Not chunking-specific; applies to plain messages too. Fix: return the permit in the drop block. Adds a test (plain messages, receiverQueueSize=1) that stalls without the fix and passes with it. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
lhotari
left a comment
There was a problem hiding this comment.
Thanks for tracking down this subtle flow-control leak and for adding a regression test with receiverQueueSize=1 that makes it visible.
I checked the premise and it holds: after a seek the consumer reconnects, ConsumerImpl.java:1071 resets the client permits and ConsumerImpl.java:975 sends a fresh full Flow, and only then does the broker re-dispatch the boundary entry. That entry consumes one broker-side permit and is dropped at ConsumerImpl.java:1541 before messageProcessed() can repay it, so the reconnect reset does not cover it. I also ran the new test with the added increaseAvailablePermits(cnx) removed and it fails (receive returns null); with the fix it passes. I found no double-count with chunked messages (non-last chunks are credited in processMessageChunk, the last chunk is repaid by the new line), zero-queue consumers, the payload-processor path or MultiTopicsConsumerImpl, and partially dropped batches are already refunded through the skipped-messages accounting. Two small suggestions inline.
Fixes #26727
Motivation
When a consumer seeks to (or is created with) a specific startMessageId, the broker re-dispatches the boundary message, which the client filters out in messageReceived() via the
isSameEntry(msgId) && isPriorEntryIndex(...)block. That block released the payload and returned without calling increaseAvailablePermits(), so the dropped message's outstanding flow-control permit was never repaid (it is not delivered, so messageProcessed() never runs for it).Each seek-to-messageId therefore leaks one permit. It is masked at normal receiver queue sizes (seek re-establishes flow control), but with receiverQueueSize=1 the single leaked permit exhausts the whole budget and the consumer stalls immediately after the seek -- the message after the seek target is never delivered. Not chunking-specific; applies to plain messages too.
Modifications
return the permit in the drop block. Adds a test (plain messages, receiverQueueSize=1) that stalls without the fix and passes with it.
Verifying this change
(Please pick either of the following options)
This change added tests and can be verified as follows:
Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes